diff --git a/Makefile b/Makefile index 567a871b04..c50e90eb40 100644 --- a/Makefile +++ b/Makefile @@ -323,6 +323,13 @@ test-upgrade: $(JUNITREPORT) IMAGE_LOGGING_EVENTROUTER=$(IMAGE_LOGGING_EVENTROUTER) \ exit 0 +.PHONY: test-e2e-from-ci-bundle +test-e2e-from-ci-bundle: $(JUNITREPORT) + RELATED_IMAGE_VECTOR=$(IMAGE_LOGGING_VECTOR) \ + RELATED_IMAGE_LOG_FILE_METRIC_EXPORTER=$(IMAGE_LOGFILEMETRICEXPORTER) \ + IMAGE_LOGGING_EVENTROUTER=$(IMAGE_LOGGING_EVENTROUTER) \ + LOG_LEVEL=3 hack/test-e2e-from-ci-bundle.sh + .PHONY: test-e2e test-e2e: $(JUNITREPORT) RELATED_IMAGE_VECTOR=$(IMAGE_LOGGING_VECTOR) \ diff --git a/hack/manifests/clusterlogforwarder/base/clusterlogforwarder.yaml b/hack/manifests/clusterlogforwarder/base/clusterlogforwarder.yaml index effbb8fedf..e3a6d0fccf 100644 --- a/hack/manifests/clusterlogforwarder/base/clusterlogforwarder.yaml +++ b/hack/manifests/clusterlogforwarder/base/clusterlogforwarder.yaml @@ -2,7 +2,6 @@ apiVersion: "observability.openshift.io/v1" kind: ClusterLogForwarder metadata: name: collector - namespace: openshift-logging spec: collector: resources: @@ -14,6 +13,6 @@ spec: - name: infra inputRefs: - application - - infrastructure +# - infrastructure outputRefs: - replaceme diff --git a/hack/manifests/clusterlogforwarder/overlays/syslog/udp/kustomization.yaml b/hack/manifests/clusterlogforwarder/overlays/syslog/udp/kustomization.yaml new file mode 100644 index 0000000000..db6aa37b5b --- /dev/null +++ b/hack/manifests/clusterlogforwarder/overlays/syslog/udp/kustomization.yaml @@ -0,0 +1,27 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +resources: +- ../../../base + +patches: +- patch: |- + - op: replace + path: /spec/pipelines/0/outputRefs/0 + value: rsyslog + - op: add + path: /spec/outputs + value: + - name: rsyslog + type: syslog + syslog: + tuning: + deliveryMode: AtLeastOnce + url: 'udp://rsyslog-server-service.openshift-logging.svc.cluster.local:514' + rfc: RFC5424 + facility: local0 + enrichment: KubernetesMinimal + appName: '{.systemd.u.SYSLOG_IDENTIFIER||.log_type||"-"}' + procId: '{.systemd.t.PID||"-"}' + msgId: '{.systemd.u.MESSAGE_ID||"-"}' + target: + kind: ClusterLogForwarder \ No newline at end of file diff --git a/hack/manifests/log-generator/base/kustomization.yaml b/hack/manifests/log-generator/base/kustomization.yaml new file mode 100644 index 0000000000..a828a12d99 --- /dev/null +++ b/hack/manifests/log-generator/base/kustomization.yaml @@ -0,0 +1,13 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +labels: + - pairs: + app.kubernetes.io/name: log-generator + app.kubernetes.io/component: log-generator + app.kubernetes.io/part-of: cluster-logging + includeSelectors: true + includeTemplates: true +resources: +- log_generator_deployment.yaml +- servicemonitor.yaml +- service.yaml diff --git a/hack/manifests/log-generator/base/log_generator_deployment.yaml b/hack/manifests/log-generator/base/log_generator_deployment.yaml new file mode 100644 index 0000000000..38a2a3e532 --- /dev/null +++ b/hack/manifests/log-generator/base/log_generator_deployment.yaml @@ -0,0 +1,35 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: log-generator +spec: + replicas: 1 + template: + spec: + securityContext: + runAsNonRoot: true + seccompProfile: + type: RuntimeDefault + containers: + - name: log-generator + env: + - name: LOGS_PER_SECOND + value: "1" + - name: PAYLOAD_SIZE + value: "512" + image: quay.io/openshift-logging/cluster-logging-load-client:latest + imagePullPolicy: Always + args: + - --logs-per-second=$(LOGS_PER_SECOND) + - --synthetic-payload-size=$(PAYLOAD_SIZE) + - --log-type=synthetic + ports: + - containerPort: 8081 + name: metrics + protocol: TCP + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: + - ALL diff --git a/hack/manifests/log-generator/base/service.yaml b/hack/manifests/log-generator/base/service.yaml new file mode 100644 index 0000000000..f6d4dd498f --- /dev/null +++ b/hack/manifests/log-generator/base/service.yaml @@ -0,0 +1,10 @@ +apiVersion: v1 +kind: Service +metadata: + name: log-generator +spec: + ports: + - name: metrics + port: 8081 + protocol: TCP + targetPort: metrics \ No newline at end of file diff --git a/hack/manifests/log-generator/base/servicemonitor.yaml b/hack/manifests/log-generator/base/servicemonitor.yaml new file mode 100644 index 0000000000..a1bd379548 --- /dev/null +++ b/hack/manifests/log-generator/base/servicemonitor.yaml @@ -0,0 +1,16 @@ +apiVersion: monitoring.coreos.com/v1 +kind: ServiceMonitor +metadata: + name: log-generator +spec: + selector: + matchLabels: + app.kubernetes.io/name: log-generator + app.kubernetes.io/component: log-generator + app.kubernetes.io/part-of: cluster-logging + endpoints: + - port: metrics + jobLabel: app.kubernetes.io/name + podTargetLabels: + - app.kubernetes.io/name + - app.kubernetes.io/component diff --git a/hack/manifests/log-generator/overlays/by_payload_size/1024b_0001lps/kustomization.yaml b/hack/manifests/log-generator/overlays/by_payload_size/1024b_0001lps/kustomization.yaml new file mode 100644 index 0000000000..c342e0042a --- /dev/null +++ b/hack/manifests/log-generator/overlays/by_payload_size/1024b_0001lps/kustomization.yaml @@ -0,0 +1,16 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +resources: +- ../../../base + +patches: +- patch: |- + - op: replace + path: /spec/template/spec/containers/0/env + value: + - name: LOGS_PER_SECOND + value: "1" + - name: PAYLOAD_SIZE + value: "1024" + target: + kind: Deployment \ No newline at end of file diff --git a/hack/manifests/log-generator/overlays/logtype/application/kustomization.yaml b/hack/manifests/log-generator/overlays/logtype/application/kustomization.yaml new file mode 100644 index 0000000000..eed8d1327f --- /dev/null +++ b/hack/manifests/log-generator/overlays/logtype/application/kustomization.yaml @@ -0,0 +1,15 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +resources: +- ../../../base + +patches: +- patch: |- + - op: replace + path: /spec/template/spec/containers/0/args + value: + - --logs-per-second=$(LOGS_PER_SECOND) + - --log-type=application + - --log-format=raw + target: + kind: Deployment \ No newline at end of file diff --git a/hack/test-e2e-from-ci-bundle.sh b/hack/test-e2e-from-ci-bundle.sh new file mode 100755 index 0000000000..ba6f33e57d --- /dev/null +++ b/hack/test-e2e-from-ci-bundle.sh @@ -0,0 +1,7 @@ +#!/bin/bash + +current_dir=$(dirname "${BASH_SOURCE[0]}" ) +GOFLAGS=-mod=mod NO_COLOR=1 go test -v -timeout=90m -ginkgo.v -ginkgo.trace \ + -ginkgo.poll-progress-after=300s \ + -ginkgo.poll-progress-interval=30s \ + "${current_dir}/../test/e2e/..." diff --git a/internal/factory/deployment.go b/internal/factory/deployment.go index 8ce8878de0..d8983ac30a 100644 --- a/internal/factory/deployment.go +++ b/internal/factory/deployment.go @@ -17,7 +17,7 @@ func NewDeployment(namespace, deploymentName, component, impl string, replicas i dpl := runtime.NewDeployment(namespace, deploymentName, visitors...) runtime.NewDeploymentBuilder(dpl).WithTemplateAnnotations(annotations). - WithTemplateLabels(dpl.Labels). + WithTemplateLabels(selectors). WithSelector(selectors). WithPodSpec(podSpec). WithReplicas(utils.GetPtr(replicas)) diff --git a/test/e2e/logforwarding/syslog/recovery_test.go b/test/e2e/logforwarding/syslog/recovery_test.go new file mode 100644 index 0000000000..5d0328080c --- /dev/null +++ b/test/e2e/logforwarding/syslog/recovery_test.go @@ -0,0 +1,142 @@ +package syslog + +import ( + "context" + "fmt" + "time" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + obs "github.com/openshift/cluster-logging-operator/api/observability/v1" + "github.com/openshift/cluster-logging-operator/internal/constants" + "github.com/openshift/cluster-logging-operator/internal/runtime" + obsruntime "github.com/openshift/cluster-logging-operator/internal/runtime/observability" + "github.com/openshift/cluster-logging-operator/internal/utils" + framework "github.com/openshift/cluster-logging-operator/test/framework/e2e" + "github.com/openshift/cluster-logging-operator/test/helpers/rand" + helpersyslog "github.com/openshift/cluster-logging-operator/test/helpers/syslog" + testruntime "github.com/openshift/cluster-logging-operator/test/runtime" + testruntimeobs "github.com/openshift/cluster-logging-operator/test/runtime/observability" + apps "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +var _ = Describe("[ClusterLogForwarder] Syslog UDP connection recovery", func() { + var ( + err error + e2e *framework.E2ETestFramework + forwarder *obs.ClusterLogForwarder + forwarderName = "my-forwarder" + deployNS string + syslogDeployment *apps.Deployment + serviceAccount *corev1.ServiceAccount + ) + + Describe("with vector collector over UDP", func() { + BeforeEach(func() { + e2e = framework.NewE2ETestFramework() + deployNS = e2e.Test.NS.Name + + By("Create log generator first so it starts producing logs") + msg := rand.Word(1024) + logGenerator := testruntime.NewLogGeneratorDeployment(deployNS, "log-generator", int32(60), 1*time.Second, string(msg)) + Expect(e2e.Test.Create(logGenerator)).To(Succeed(), "failed to create log generator") + + if serviceAccount, err = e2e.BuildAuthorizationFor(deployNS, forwarderName). + AllowClusterRole(framework.ClusterRoleCollectApplicationLogs). + AllowClusterRole(framework.ClusterRoleCollectInfrastructureLogs). + AllowClusterRole(framework.ClusterRoleCollectAuditLogs).Create(); err != nil { + Fail(err.Error()) + } + + By("Deploy syslog UDP receiver first") + syslogDeployment, err = e2e.DeploySyslogReceiver(deployNS, corev1.ProtocolUDP, false, helpersyslog.RFC5424) + Expect(err).To(BeNil(), "should successfully deploy syslog UDP receiver") + + By("Create forwarder with the syslog receiver running") + forwarder = testruntimeobs.NewClusterLogForwarderBuilder(obsruntime.NewClusterLogForwarder(deployNS, forwarderName, runtime.Initialize), func(clf *obs.ClusterLogForwarder) { + clf.Spec.Collector = &obs.CollectorSpec{ + Resources: &corev1.ResourceRequirements{ + Limits: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("1"), + corev1.ResourceMemory: resource.MustParse("2Gi"), + }, + Requests: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("100m"), + corev1.ResourceMemory: resource.MustParse("64Mi"), + }, + }, + } + clf.Spec.ServiceAccount.Name = serviceAccount.Name + }). + FromInput(obs.InputTypeApplication). + ToSyslogOutput(obs.SyslogRFC5424, func(output *obs.OutputSpec) { + output.Syslog = &obs.Syslog{ + URL: fmt.Sprintf("udp://%s.%s.svc:514", framework.SyslogReceiverName, deployNS), + RFC: obs.SyslogRFC5424, + Tuning: &obs.SyslogTuningSpec{ + DeliveryMode: obs.DeliveryModeAtLeastOnce, + }, + Facility: "local0", + Enrichment: obs.EnrichmentTypeKubernetesMinimal, + AppName: `{.systemd.u.SYSLOG_IDENTIFIER||.log_type||"-"}`, + ProcId: `{.systemd.t.PID||"-"}`, + MsgId: `{.systemd.u.MESSAGE_ID||"-"}`, + } + }).End() + + if err := e2e.CreateObservabilityClusterLogForwarder(forwarder); err != nil { + Fail(fmt.Sprintf("Unable to create an instance of logforwarder: %v", err)) + } + + By("Waiting for the collector daemonset") + if err := e2e.WaitForDaemonSet(forwarder.Namespace, forwarder.Name); err != nil { + Fail(err.Error()) + } + }) + + Context("should recover and send logs after syslog UDP receiver restarts", func() { + var podDeleteTest = func() { + By("Verify initial logs are flowing") + logStore := e2e.LogStores[syslogDeployment.GetName()] + Expect(logStore.HasApplicationLogs(time.Minute*2)).To(BeTrue(), "expected to collect application logs initially") + + ctx := context.TODO() + labelSelector := fmt.Sprintf("%s=%s", constants.LabelK8sComponent, syslogDeployment.Name) + + By("Delete syslog receiver pod repeatedly to trigger ICMP error caching on vector's UDP socket") + for i := 0; i < 10; i++ { + err = e2e.KubeClient.CoreV1().Pods(deployNS).DeleteCollection(ctx, + metav1.DeleteOptions{GracePeriodSeconds: utils.GetPtr[int64](0)}, + metav1.ListOptions{LabelSelector: labelSelector}, + ) + Expect(err).To(BeNil(), "should be able to delete syslog receiver pods") + time.Sleep(5 * time.Second) + } + Expect(e2e.WaitForDeployment(deployNS, syslogDeployment.Name, time.Second, time.Second*30)).To(Succeed(), "replicas should become available ") + + By("Verify that logs resume flowing after recovery") + Expect(logStore.HasApplicationLogs(1*time.Minute)).To(BeTrue(), "expected to collect application logs after receiver recovers") + //Fail("I should never get here") + } + + It("when using a ClusterIP service", func() { + podDeleteTest() + }) + + It("when using a NodePort service", func() { + Expect(e2e.Test.Delete(runtime.NewService(deployNS, framework.SyslogReceiverName))).To(Succeed(), "should delete syslog receiver service") + By("Replacing the ClusterIP service with a NodePort service") + _, err = e2e.CreateSyslogService(syslogDeployment, corev1.ProtocolUDP, corev1.ServiceTypeNodePort) + Expect(err).To(BeNil(), "should successfully create NodePort service") + podDeleteTest() + }) + }) + + AfterEach(func() { + e2e.Cleanup() + }) + }) +}) diff --git a/test/e2e/logforwarding/syslog/syslog_suite_test.go b/test/e2e/logforwarding/syslog/syslog_suite_test.go new file mode 100644 index 0000000000..e5294a0ad1 --- /dev/null +++ b/test/e2e/logforwarding/syslog/syslog_suite_test.go @@ -0,0 +1,13 @@ +package syslog + +import ( + "testing" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +func TestSyslog(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "Syslog Log Forwarding E2E Suite") +} diff --git a/test/framework/e2e/framework.go b/test/framework/e2e/framework.go index 32a43ef43b..ee808737a4 100644 --- a/test/framework/e2e/framework.go +++ b/test/framework/e2e/framework.go @@ -237,7 +237,7 @@ func (tc *E2ETestFramework) Create(obj crclient.Object) error { tc.AddCleanup(func() error { return tc.Test.Delete(obj) }) - clolog.Info("Creating object", "obj", string(body)) + clolog.V(2).Info("Creating object", "obj", string(body)) return tc.Test.Recreate(obj) } @@ -247,7 +247,7 @@ func (tc *E2ETestFramework) CreateObservabilityClusterLogForwarder(forwarder *ob func DoCleanup() bool { doCleanup := strings.TrimSpace(os.Getenv("DO_CLEANUP")) - clolog.Info("Running Cleanup script ....", "DO_CLEANUP", doCleanup) + clolog.V(1).Info("Running Cleanup script ....", "DO_CLEANUP", doCleanup) return doCleanup == "" || strings.ToLower(doCleanup) == "true" } @@ -264,7 +264,7 @@ func (tc *E2ETestFramework) Cleanup() { } else { clolog.V(1).Info("Test passed. Skipping artifacts gathering") } - clolog.Info("Running e2e cleanup functions, ", "number", len(tc.CleanupFns)) + clolog.V(1).Info("Running e2e cleanup functions, ", "number", len(tc.CleanupFns)) for _, cleanup := range tc.CleanupFns { clolog.V(5).Info("Running an e2e cleanup function") if err := cleanup(); err != nil { @@ -282,14 +282,14 @@ func RunCleanupScript() { clolog.Info("No cleanup script provided") return } - clolog.Info("Script", "CLEANUP_CMD", value) + clolog.V(1).Info("Script", "CLEANUP_CMD", value) args := strings.Split(value, " ") // #nosec G204 cmd := exec.Command(args[0], args[1:]...) cmd.Env = nil result, err := cmd.CombinedOutput() - clolog.Info("RunCleanupScript output: ", "output", string(result)) - clolog.Info("RunCleanupScript err: ", "error", err) + clolog.V(2).Info("RunCleanupScript output: ", "output", string(result)) + clolog.V(2).Info("RunCleanupScript err: ", "error", err) } } diff --git a/test/framework/e2e/syslog.go b/test/framework/e2e/syslog.go index d7eac47b0c..949edd651b 100644 --- a/test/framework/e2e/syslog.go +++ b/test/framework/e2e/syslog.go @@ -4,15 +4,16 @@ import ( "context" "errors" "fmt" - "log" "strconv" "strings" "time" + "github.com/openshift/cluster-logging-operator/internal/constants" "github.com/openshift/cluster-logging-operator/test/helpers/certificate" + "github.com/openshift/cluster-logging-operator/test/helpers/syslog" rbacv1 "k8s.io/api/rbac/v1" + "k8s.io/apimachinery/pkg/util/intstr" - "github.com/openshift/cluster-logging-operator/internal/constants" "github.com/openshift/cluster-logging-operator/internal/factory" "github.com/openshift/cluster-logging-operator/internal/runtime" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -35,64 +36,14 @@ type syslogReceiverLogStore struct { const ( SyslogReceiverName = "syslog-receiver" - ImageRemoteSyslog = "registry.redhat.io/rhel8/rsyslog:8.7-9" ) -// SyslogRfc type is the rfc used for sending syslog -type SyslogRfc int - -const ( - // RFC3164 rfc3164 - RFC3164 SyslogRfc = iota - // RFC5424 rfc5424 - RFC5424 - // RFC3164RFC5424 either rfc3164 or rfc5424 - RFC3164RFC5424 -) - -func MustParseRFC(rfc string) SyslogRfc { - switch strings.ToUpper(rfc) { - case "RFC3164": - return RFC3164 - case "RFC5424": - return RFC5424 - case "RFC3164 or RFC5424": - return RFC3164RFC5424 - } - log.Fatal("Unable to parse RFC", "rfc", rfc) - return 0 -} - -func (e SyslogRfc) String() string { - switch e { - case RFC3164: - return "RFC3164" - case RFC5424: - return "RFC5424" - case RFC3164RFC5424: - return "RFC3164 or RFC5424" - default: - return "Unknown rfc" - } -} - -func GenerateRsyslogConf(conf string, rfc SyslogRfc) string { - switch rfc { - case RFC5424: - return strings.Join([]string{conf, RuleSetRfc5424}, "\n") - case RFC3164: - return strings.Join([]string{conf, RuleSetRfc3164}, "\n") - case RFC3164RFC5424: - return strings.Join([]string{conf, RuleSetRfc3164Rfc5424}, "\n") - } - return "Invalid Conf" -} - -func (syslog *syslogReceiverLogStore) hasLogs(file string, timeToWait time.Duration) (bool, error) { +func (s *syslogReceiverLogStore) hasLogs(file string, timeToWait time.Duration) (bool, error) { options := metav1.ListOptions{ - LabelSelector: "component=syslog-receiver", + LabelSelector: constants.LabelK8sComponent + "=" + s.deployment.Name, } - pods, err := syslog.tc.KubeClient.CoreV1().Pods(constants.OpenshiftNS).List(context.TODO(), options) + clolog.V(3).Info("Listing syslog pods", "namespace", s.deployment.Namespace, "options", options) + pods, err := s.tc.KubeClient.CoreV1().Pods(s.deployment.Namespace).List(context.TODO(), options) if err != nil { return false, err } @@ -101,12 +52,14 @@ func (syslog *syslogReceiverLogStore) hasLogs(file string, timeToWait time.Durat } podName := pods.Items[0].Name cmd := fmt.Sprintf("ls %s | wc -l", file) - err = wait.PollUntilContextTimeout(context.TODO(), defaultRetryInterval, timeToWait, true, func(cxt context.Context) (done bool, err error) { - output, err := syslog.tc.PodExec(constants.OpenshiftNS, podName, "syslog-receiver", []string{"bash", "-c", cmd}) + clolog.V(3).Info("pod exec", "pod", podName, "cmd", cmd) + err = wait.PollUntilContextTimeout(context.TODO(), 1*time.Second, timeToWait, true, func(cxt context.Context) (done bool, err error) { + output, err := s.tc.PodExec(s.deployment.Namespace, podName, "syslog-receiver", []string{"bash", "-c", cmd}) if err != nil { clolog.Error(err, "failed to fetch logs from syslog-receiver") return false, nil } + clolog.V(3).Info("syslog-receiver pod exec", "pod", podName, "output", output) value, err := strconv.Atoi(strings.TrimSpace(output)) if err != nil { clolog.V(2).Error(err, "Error parsing output", "output", output) @@ -120,12 +73,12 @@ func (syslog *syslogReceiverLogStore) hasLogs(file string, timeToWait time.Durat return true, err } -func (syslog *syslogReceiverLogStore) grepLogs(expr string, logfile string, timeToWait time.Duration) (string, error) { +func (s *syslogReceiverLogStore) grepLogs(expr string, logfile string, timeToWait time.Duration) (string, error) { NotFound := "No Found" options := metav1.ListOptions{ - LabelSelector: "component=syslog-receiver", + LabelSelector: constants.LabelK8sComponent + "=" + s.deployment.Name, } - pods, err := syslog.tc.KubeClient.CoreV1().Pods(constants.OpenshiftNS).List(context.TODO(), options) + pods, err := s.tc.KubeClient.CoreV1().Pods(s.deployment.Namespace).List(context.TODO(), options) if err != nil { return NotFound, err } @@ -137,8 +90,8 @@ func (syslog *syslogReceiverLogStore) grepLogs(expr string, logfile string, time clolog.V(3).Info("running expression", "expression", cmd) var value string - err = wait.PollUntilContextTimeout(context.TODO(), defaultRetryInterval, timeToWait, true, func(cxt context.Context) (done bool, err error) { - output, err := syslog.tc.PodExec(constants.OpenshiftNS, pods.Items[0].Name, "syslog-receiver", []string{"bash", "-c", cmd}) + err = wait.PollUntilContextTimeout(context.TODO(), 1*time.Second, timeToWait, true, func(cxt context.Context) (done bool, err error) { + output, err := s.tc.PodExec(s.deployment.Namespace, pods.Items[0].Name, "syslog-receiver", []string{"bash", "-c", cmd}) if err != nil { clolog.Error(err, "failed to fetch logs from syslog-receiver") return false, nil @@ -152,51 +105,34 @@ func (syslog *syslogReceiverLogStore) grepLogs(expr string, logfile string, time return value, nil } -func (syslog *syslogReceiverLogStore) ApplicationLogs(timeToWait time.Duration) (types.Logs, error) { +func (s *syslogReceiverLogStore) ApplicationLogs(timeToWait time.Duration) (types.Logs, error) { panic("Method not implemented") } -func (syslog *syslogReceiverLogStore) HasInfraStructureLogs(timeToWait time.Duration) (bool, error) { - return syslog.hasLogs("/tmp/infra.log", timeToWait) +func (s *syslogReceiverLogStore) HasInfraStructureLogs(timeToWait time.Duration) (bool, error) { + return s.hasLogs("/tmp/infra.log", timeToWait) } -func (syslog *syslogReceiverLogStore) HasApplicationLogs(timeToWait time.Duration) (bool, error) { - return false, fmt.Errorf("not implemented") +func (s *syslogReceiverLogStore) HasApplicationLogs(timeToWait time.Duration) (bool, error) { + return s.hasLogs("/tmp/app.log", timeToWait) } -func (syslog *syslogReceiverLogStore) HasAuditLogs(timeToWait time.Duration) (bool, error) { +func (s *syslogReceiverLogStore) HasAuditLogs(timeToWait time.Duration) (bool, error) { return false, fmt.Errorf("not implemented") } -func (syslog *syslogReceiverLogStore) GrepLogs(expr string, timeToWait time.Duration) (string, error) { - return syslog.grepLogs(expr, "/tmp/infra.log", timeToWait) +func (s *syslogReceiverLogStore) GrepLogs(expr string, timeToWait time.Duration) (string, error) { + return s.grepLogs(expr, "/tmp/infra.log", timeToWait) } -func (syslog *syslogReceiverLogStore) RetrieveLogs() (map[string]string, error) { +func (s *syslogReceiverLogStore) RetrieveLogs() (map[string]string, error) { return nil, fmt.Errorf("not implemented") } -func (syslog *syslogReceiverLogStore) ClusterLocalEndpoint() string { +func (s *syslogReceiverLogStore) ClusterLocalEndpoint() string { panic("not implemented") } -func (tc *E2ETestFramework) createSyslogServiceAccount() (serviceAccount *corev1.ServiceAccount, err error) { - opts := metav1.CreateOptions{} - serviceAccount = runtime.NewServiceAccount(constants.OpenshiftNS, "syslog-receiver") - if serviceAccount, err = tc.KubeClient.CoreV1().ServiceAccounts(constants.OpenshiftNS).Create(context.TODO(), serviceAccount, opts); err != nil { - return nil, err - } - tc.AddCleanup(func() error { - opts := metav1.DeleteOptions{} - err := tc.KubeClient.CoreV1().ServiceAccounts(constants.OpenshiftNS).Delete(context.TODO(), serviceAccount.Name, opts) - if apierrors.IsNotFound(err) { - return nil - } - return err - }) - return serviceAccount, nil -} - func (tc *E2ETestFramework) CreateLegacySyslogConfigMap(namespace, conf string) (err error) { opts := metav1.CreateOptions{} fluentdConfigMap := runtime.NewConfigMap( @@ -221,10 +157,10 @@ func (tc *E2ETestFramework) CreateLegacySyslogConfigMap(namespace, conf string) return nil } -func (tc *E2ETestFramework) createSyslogRbac(name string) (err error) { +func (tc *E2ETestFramework) createSyslogRbac(namespace, name string) (err error) { opts := metav1.CreateOptions{} saRole := runtime.NewRole( - constants.OpenshiftNS, + namespace, name, runtime.NewPolicyRules( runtime.NewPolicyRule( @@ -236,13 +172,13 @@ func (tc *E2ETestFramework) createSyslogRbac(name string) (err error) { )..., ) - if _, err = tc.KubeClient.RbacV1().Roles(constants.OpenshiftNS).Create(context.TODO(), saRole, opts); err != nil { + if _, err = tc.KubeClient.RbacV1().Roles(namespace).Create(context.TODO(), saRole, opts); err != nil { return err } tc.AddCleanup(func() error { opts := metav1.DeleteOptions{} - err := tc.KubeClient.RbacV1().Roles(constants.OpenshiftNS).Delete(context.TODO(), name, opts) + err := tc.KubeClient.RbacV1().Roles(namespace).Delete(context.TODO(), name, opts) if apierrors.IsNotFound(err) { return nil } @@ -258,7 +194,7 @@ func (tc *E2ETestFramework) createSyslogRbac(name string) (err error) { subject.APIGroup = "" roleBinding := runtime.NewRoleBinding( - constants.OpenshiftNS, + namespace, name, rbacv1.RoleRef{ Kind: "Role", @@ -270,12 +206,12 @@ func (tc *E2ETestFramework) createSyslogRbac(name string) (err error) { )..., ) - if _, err = tc.KubeClient.RbacV1().RoleBindings(constants.OpenshiftNS).Create(context.TODO(), roleBinding, rbOpts); err != nil { + if _, err = tc.KubeClient.RbacV1().RoleBindings(namespace).Create(context.TODO(), roleBinding, rbOpts); err != nil { return err } tc.AddCleanup(func() error { opts := metav1.DeleteOptions{} - err := tc.KubeClient.RbacV1().RoleBindings(constants.OpenshiftNS).Delete(context.TODO(), name, opts) + err := tc.KubeClient.RbacV1().RoleBindings(namespace).Delete(context.TODO(), name, opts) if apierrors.IsNotFound(err) { return nil } @@ -284,27 +220,27 @@ func (tc *E2ETestFramework) createSyslogRbac(name string) (err error) { return nil } -func (tc *E2ETestFramework) DeploySyslogReceiver(testDir string, protocol corev1.Protocol, withTLS bool, rfc SyslogRfc) (deployment *apps.Deployment, err error) { +func (tc *E2ETestFramework) DeploySyslogReceiver(namespace string, protocol corev1.Protocol, withTLS bool, rfc syslog.SyslogRfc) (deployment *apps.Deployment, err error) { logStore := &syslogReceiverLogStore{ tc: tc, } - serviceAccount, err := tc.createSyslogServiceAccount() + serviceAccount, err := tc.createServiceAccount(namespace, SyslogReceiverName) if err != nil { return nil, err } - if err := tc.createSyslogRbac(SyslogReceiverName); err != nil { + if err := tc.createSyslogRbac(namespace, SyslogReceiverName); err != nil { return nil, err } container := corev1.Container{ Name: SyslogReceiverName, - Image: ImageRemoteSyslog, + Image: syslog.ImageRemoteSyslog, ImagePullPolicy: corev1.PullAlways, - Args: []string{"rsyslogd", "-n", "-f", "/rsyslog/etc/rsyslog.conf"}, + Command: []string{"/usr/sbin/rsyslogd", "-i", "/tmp/rsyslog.pid", "-n", "-f", "/etc/rsyslog/rsyslog.conf"}, VolumeMounts: []corev1.VolumeMount{ { Name: "config", ReadOnly: true, - MountPath: "/rsyslog/etc", + MountPath: "/etc/rsyslog", }, }, SecurityContext: &corev1.SecurityContext{ @@ -326,6 +262,12 @@ func (tc *E2ETestFramework) DeploySyslogReceiver(testDir string, protocol corev1 }, }, }, + { + Name: "log", + VolumeSource: corev1.VolumeSource{ + EmptyDir: &corev1.EmptyDirVolumeSource{}, + }, + }, }, ServiceAccountName: serviceAccount.Name, SecurityContext: &corev1.PodSecurityContext{ @@ -339,27 +281,27 @@ func (tc *E2ETestFramework) DeploySyslogReceiver(testDir string, protocol corev1 var rsyslogConf string switch protocol { case corev1.ProtocolUDP: - rsyslogConf = UdpSyslogInput + rsyslogConf = syslog.UdpSyslogInput default: - rsyslogConf = TcpSyslogInput + rsyslogConf = syslog.TcpSyslogInput } if withTLS { switch protocol { case corev1.ProtocolUDP: - rsyslogConf = UdpSyslogInputWithTLS + rsyslogConf = syslog.UdpSyslogInputWithTLS default: - rsyslogConf = TcpSyslogInputWithTLS + rsyslogConf = syslog.TcpSyslogInputWithTLS } - secret, err := tc.CreateSyslogReceiverSecrets(testDir, SyslogReceiverName, SyslogReceiverName) + secret, err := tc.CreateSyslogReceiverSecrets(namespace, SyslogReceiverName, SyslogReceiverName) if err != nil { return nil, err } tc.AddCleanup(func() error { opts := metav1.DeleteOptions{} - err := tc.KubeClient.CoreV1().Secrets(constants.OpenshiftNS).Delete(context.TODO(), SyslogReceiverName, opts) + err := tc.KubeClient.CoreV1().Secrets(namespace).Delete(context.TODO(), SyslogReceiverName, opts) if apierrors.IsNotFound(err) { return nil } @@ -380,19 +322,19 @@ func (tc *E2ETestFramework) DeploySyslogReceiver(testDir string, protocol corev1 }) } - rsyslogConf = GenerateRsyslogConf(rsyslogConf, rfc) + rsyslogConf = syslog.GenerateRsyslogConf(rsyslogConf, rfc) cOpts := metav1.CreateOptions{} - config := runtime.NewConfigMap(constants.OpenshiftNS, container.Name, map[string]string{ + config := runtime.NewConfigMap(namespace, container.Name, map[string]string{ "rsyslog.conf": rsyslogConf, }) - config, err = tc.KubeClient.CoreV1().ConfigMaps(constants.OpenshiftNS).Create(context.TODO(), config, cOpts) + config, err = tc.KubeClient.CoreV1().ConfigMaps(namespace).Create(context.TODO(), config, cOpts) if err != nil { return nil, err } tc.AddCleanup(func() error { opts := metav1.DeleteOptions{} - err := tc.KubeClient.CoreV1().ConfigMaps(constants.OpenshiftNS).Delete(context.TODO(), config.Name, opts) + err := tc.KubeClient.CoreV1().ConfigMaps(namespace).Delete(context.TODO(), config.Name, opts) if apierrors.IsNotFound(err) { return nil } @@ -401,69 +343,87 @@ func (tc *E2ETestFramework) DeploySyslogReceiver(testDir string, protocol corev1 dOpts := metav1.CreateOptions{} syslogDeployment := factory.NewDeployment( - constants.OpenshiftNS, + namespace, container.Name, container.Name, serviceAccount.Name, 1, podSpec, ) - - // Add instance label to pod spec template. Service now selects using instance name as well - syslogDeployment.Spec.Template.Labels[constants.LabelK8sInstance] = serviceAccount.Name - - syslogDeployment, err = tc.KubeClient.AppsV1().Deployments(constants.OpenshiftNS).Create(context.TODO(), syslogDeployment, dOpts) + syslogDeployment, err = tc.KubeClient.AppsV1().Deployments(namespace).Create(context.TODO(), syslogDeployment, dOpts) if err != nil { return nil, err } - service := factory.NewService( - serviceAccount.Name, - constants.OpenshiftNS, - serviceAccount.Name, - serviceAccount.Name, - []corev1.ServicePort{ - { - Protocol: protocol, - Port: 24224, - }, - }, - ) + + if _, err = tc.CreateSyslogService(syslogDeployment, protocol, corev1.ServiceTypeClusterIP); err != nil { + return nil, err + } tc.AddCleanup(func() error { var zerograce int64 deleteopts := metav1.DeleteOptions{ GracePeriodSeconds: &zerograce, } - err := tc.KubeClient.AppsV1().Deployments(constants.OpenshiftNS).Delete(context.TODO(), syslogDeployment.Name, deleteopts) + err := tc.KubeClient.AppsV1().Deployments(namespace).Delete(context.TODO(), syslogDeployment.Name, deleteopts) if apierrors.IsNotFound(err) { return nil } return err }) + logStore.deployment = syslogDeployment + + name := syslogDeployment.GetName() + tc.LogStores[name] = logStore + return syslogDeployment, tc.WaitForDeployment(namespace, syslogDeployment.Name, defaultRetryInterval, defaultTimeout) +} + +func (tc *E2ETestFramework) CreateSyslogService(syslogDeployment *apps.Deployment, protocol corev1.Protocol, serviceType corev1.ServiceType) (service *corev1.Service, err error) { + + service = factory.NewService( + syslogDeployment.Name, + syslogDeployment.Namespace, + syslogDeployment.Name, + syslogDeployment.Name, + []corev1.ServicePort{ + { + Name: "udp", + Protocol: protocol, + TargetPort: intstr.FromInt32(24224), + Port: 514, + }, + { + Name: "tcp", + Protocol: corev1.ProtocolTCP, + TargetPort: intstr.FromInt32(24224), + Port: 514, + }, + }, + func(o runtime.Object) { + runtime.SetCommonLabels(o, syslogDeployment.Name, syslogDeployment.Name, syslogDeployment.Name) + }, + ) + service.Spec.Type = serviceType + sOpts := metav1.CreateOptions{} - service, err = tc.KubeClient.CoreV1().Services(constants.OpenshiftNS).Create(context.TODO(), service, sOpts) + service, err = tc.KubeClient.CoreV1().Services(syslogDeployment.Namespace).Create(context.TODO(), service, sOpts) if err != nil { return nil, err } tc.AddCleanup(func() error { opts := metav1.DeleteOptions{} - err := tc.KubeClient.CoreV1().Services(constants.OpenshiftNS).Delete(context.TODO(), service.Name, opts) + err = tc.KubeClient.CoreV1().Services(syslogDeployment.Namespace).Delete(context.TODO(), service.Name, opts) if apierrors.IsNotFound(err) { return nil } return err }) - logStore.deployment = syslogDeployment - - name := syslogDeployment.GetName() - tc.LogStores[name] = logStore - return syslogDeployment, tc.WaitForDeployment(constants.OpenshiftNS, syslogDeployment.Name, defaultRetryInterval, defaultTimeout) + return service, nil } -func (tc *E2ETestFramework) CreateSyslogReceiverSecrets(testDir, logStoreName, secretName string) (secret *corev1.Secret, err error) { +func (tc *E2ETestFramework) CreateSyslogReceiverSecrets(namespace, logStoreName, secretName string) (secret *corev1.Secret, err error) { ca := certificate.NewCA(nil, "Root CA") // Self-signed CA - serverCert := certificate.NewCert(ca, "", logStoreName, fmt.Sprintf("%s.%s.svc", logStoreName, constants.OpenshiftNS)) + serverCert := certificate.NewCert(ca, "", logStoreName, fmt.Sprintf("%s.%s.svc", logStoreName, namespace)) data := map[string][]byte{ "tls.key": serverCert.PrivateKeyPEM(), @@ -474,12 +434,12 @@ func (tc *E2ETestFramework) CreateSyslogReceiverSecrets(testDir, logStoreName, s sOpts := metav1.CreateOptions{} secret = runtime.NewSecret( - constants.OpenshiftNS, + namespace, secretName, data, ) clolog.V(3).Info("Creating secret for logStore", "secret", secret.Name, "logStore", logStoreName) - if secret, err = tc.KubeClient.CoreV1().Secrets(constants.OpenshiftNS).Create(context.TODO(), secret, sOpts); err != nil { + if secret, err = tc.KubeClient.CoreV1().Secrets(namespace).Create(context.TODO(), secret, sOpts); err != nil { return nil, err } return secret, nil diff --git a/test/framework/e2e/wait.go b/test/framework/e2e/wait.go index f69ace3b94..b8b4f07edb 100644 --- a/test/framework/e2e/wait.go +++ b/test/framework/e2e/wait.go @@ -38,7 +38,7 @@ func (tc *E2ETestFramework) WaitForDaemonSet(namespace, name string) error { ds, err := tc.KubeClient.AppsV1().DaemonSets(namespace).Get(context.TODO(), name, metav1.GetOptions{}) if err != nil { - clolog.V(0).Error(err, "error polling and waiting for daemonset", "namespace", namespace, "name", name) + clolog.V(3).Error(err, "error polling and waiting for daemonset", "namespace", namespace, "name", name) return false, nil } if ds.Status.DesiredNumberScheduled == ds.Status.NumberReady { diff --git a/test/framework/functional/constants.go b/test/framework/functional/constants.go index 7020ccde2f..9faa590394 100644 --- a/test/framework/functional/constants.go +++ b/test/framework/functional/constants.go @@ -48,7 +48,7 @@ var ( string(obs.InputTypeInfrastructure): ApplicationLogFile, }, string(obs.OutputTypeSyslog): { - applicationLog: "/tmp/infra.log", + applicationLog: "/tmp/app.log", auditLog: "/tmp/infra.log", k8sAuditLog: "/tmp/infra.log", ovnAuditLog: "/tmp/infra.log", diff --git a/test/framework/functional/output_syslog.go b/test/framework/functional/output_syslog.go index 7252158693..878469f558 100644 --- a/test/framework/functional/output_syslog.go +++ b/test/framework/functional/output_syslog.go @@ -3,13 +3,13 @@ package functional import ( obs "github.com/openshift/cluster-logging-operator/api/observability/v1" internalobs "github.com/openshift/cluster-logging-operator/internal/api/observability" + "github.com/openshift/cluster-logging-operator/test/helpers/syslog" + "net/url" "strings" - "github.com/openshift/cluster-logging-operator/internal/runtime" - "github.com/openshift/cluster-logging-operator/test/framework/e2e" - log "github.com/ViaQ/logerr/v2/log/static" + "github.com/openshift/cluster-logging-operator/internal/runtime" ) const IncreaseRsyslogMaxMessageSize = "$MaxMessageSize 50000" @@ -20,23 +20,23 @@ func (f *CollectorFunctionalFramework) AddSyslogOutput(b *runtime.PodBuilder, ou var baseRsyslogConfig string u, _ := url.Parse(output.Syslog.URL) if strings.ToLower(u.Scheme) == "udp" { - baseRsyslogConfig = e2e.UdpSyslogInput + baseRsyslogConfig = syslog.UdpSyslogInput if output.TLS != nil { - baseRsyslogConfig = e2e.UdpSyslogInputWithTLS + baseRsyslogConfig = syslog.UdpSyslogInputWithTLS } } else { - baseRsyslogConfig = e2e.TcpSyslogInput + baseRsyslogConfig = syslog.TcpSyslogInput if output.TLS != nil { - baseRsyslogConfig = e2e.TcpSyslogInputWithTLS + baseRsyslogConfig = syslog.TcpSyslogInputWithTLS } } // using unsecure rsyslog conf - rfc := e2e.RFC5424 + rfc := syslog.RFC5424 if output.Syslog != nil && output.Syslog.RFC != "" { - rfc = e2e.MustParseRFC(string(output.Syslog.RFC)) + rfc = syslog.MustParseRFC(string(output.Syslog.RFC)) } - rsyslogConf := e2e.GenerateRsyslogConf(baseRsyslogConfig, rfc) + rsyslogConf := syslog.GenerateRsyslogConf(baseRsyslogConfig, rfc) rsyslogConf = strings.Join([]string{IncreaseRsyslogMaxMessageSize, rsyslogConf}, "\n") config := runtime.NewConfigMap(b.Pod.Namespace, name, map[string]string{ "rsyslog.conf": rsyslogConf, @@ -47,7 +47,7 @@ func (f *CollectorFunctionalFramework) AddSyslogOutput(b *runtime.PodBuilder, ou } log.V(2).Info("Adding container", "name", name) - containerBuilder := b.AddContainer(name, e2e.ImageRemoteSyslog). + containerBuilder := b.AddContainer(name, syslog.ImageRemoteSyslog). AddVolumeMount(config.Name, "/rsyslog/etc", "", false). WithCmdArgs([]string{"rsyslogd", "-n", "-f", "/rsyslog/etc/rsyslog.conf"}). WithPrivilege() diff --git a/test/helpers/syslog/sender.go b/test/framework/functional/outputs/syslog/sender.go similarity index 91% rename from test/helpers/syslog/sender.go rename to test/framework/functional/outputs/syslog/sender.go index 4a4724f619..4020a75e54 100644 --- a/test/helpers/syslog/sender.go +++ b/test/framework/functional/outputs/syslog/sender.go @@ -27,5 +27,5 @@ func WriteToSyslogInputWithNetcat(framework *functional.CollectorFunctionalFrame return err } } - return fmt.Errorf("WriteToHttpInput: no HTTP input named %s", inputName) + return fmt.Errorf("WriteToSyslogInputWithNetcat: no syslog input named %s", inputName) } diff --git a/test/functional/inputs/syslog/syslog_input_test.go b/test/functional/inputs/syslog/syslog_input_test.go index 63a9a83ea5..7ddfaeb5ff 100644 --- a/test/functional/inputs/syslog/syslog_input_test.go +++ b/test/functional/inputs/syslog/syslog_input_test.go @@ -4,16 +4,17 @@ import ( "context" "encoding/json" "fmt" + "strings" + "time" + . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" obs "github.com/openshift/cluster-logging-operator/api/observability/v1" "github.com/openshift/cluster-logging-operator/internal/runtime" "github.com/openshift/cluster-logging-operator/test/framework/functional" - "github.com/openshift/cluster-logging-operator/test/helpers/syslog" + "github.com/openshift/cluster-logging-operator/test/framework/functional/outputs/syslog" testruntime "github.com/openshift/cluster-logging-operator/test/runtime/observability" "k8s.io/apimachinery/pkg/util/wait" - "strings" - "time" ) const ( diff --git a/test/helpers/syslog/syslog.go b/test/helpers/syslog/syslog.go new file mode 100644 index 0000000000..bc56b20e89 --- /dev/null +++ b/test/helpers/syslog/syslog.go @@ -0,0 +1,60 @@ +package syslog + +import ( + "log" + "strings" +) + +const ( + ImageRemoteSyslog = "registry.redhat.io/rhel9/rsyslog:9.8-1780596788" +) + +// SyslogRfc type is the rfc used for sending syslog +type SyslogRfc int + +const ( + // RFC3164 rfc3164 + RFC3164 SyslogRfc = iota + // RFC5424 rfc5424 + RFC5424 + // RFC3164RFC5424 either rfc3164 or rfc5424 + RFC3164RFC5424 +) + +func MustParseRFC(rfc string) SyslogRfc { + switch strings.ToUpper(rfc) { + case "RFC3164": + return RFC3164 + case "RFC5424": + return RFC5424 + case "RFC3164 OR RFC5424": + return RFC3164RFC5424 + } + log.Fatal("Unable to parse RFC", "rfc", rfc) + return 0 +} + +func (e SyslogRfc) String() string { + switch e { + case RFC3164: + return "RFC3164" + case RFC5424: + return "RFC5424" + case RFC3164RFC5424: + return "RFC3164 or RFC5424" + default: + return "Unknown rfc" + } +} + +func GenerateRsyslogConf(conf string, rfc SyslogRfc) string { + switch rfc { + case RFC5424: + return strings.Join([]string{conf, RuleSetRfc5424}, "\n") + case RFC3164: + return strings.Join([]string{conf, RuleSetRfc3164}, "\n") + case RFC3164RFC5424: + return strings.Join([]string{conf, RuleSetRfc3164Rfc5424}, "\n") + } + return "Invalid Conf" +} diff --git a/test/framework/e2e/syslog_conf.go b/test/helpers/syslog/syslog_conf.go similarity index 87% rename from test/framework/e2e/syslog_conf.go rename to test/helpers/syslog/syslog_conf.go index 67ec45fac6..141514425b 100644 --- a/test/framework/e2e/syslog_conf.go +++ b/test/helpers/syslog/syslog_conf.go @@ -1,4 +1,4 @@ -package e2e +package syslog const ( TcpSyslogInput = ` @@ -54,7 +54,12 @@ input(type="imudp" port="24224" ruleset="test") RuleSetRfc5424 = ` #### RULES #### ruleset(name="test" parser=["rsyslog.rfc5424"]){ - action(type="omfile" file="/tmp/infra.log" Template="RSYSLOG_SyslogProtocol23Format") + # Check message content for log_type field + if ($msg contains "\"log_type\":\"application\"") then { + action(type="omfile" file="/tmp/app.log" Template="RSYSLOG_SyslogProtocol23Format") + } else { + action(type="omfile" file="/tmp/infra.log" Template="RSYSLOG_SyslogProtocol23Format") + } } ` diff --git a/test/runtime/log_generator.go b/test/runtime/log_generator.go index 12a7e5f903..382f801ada 100644 --- a/test/runtime/log_generator.go +++ b/test/runtime/log_generator.go @@ -4,13 +4,42 @@ import ( "fmt" "time" + "github.com/openshift/cluster-logging-operator/internal/constants" "github.com/openshift/cluster-logging-operator/internal/utils" + appsv1 "k8s.io/api/apps/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/openshift/cluster-logging-operator/internal/runtime" corev1 "k8s.io/api/core/v1" ) +// NewLogGeneratorDeployment creates a single container log generator deployment +func NewLogGeneratorDeployment(namespace, name string, replicas int32, delay time.Duration, message string) *appsv1.Deployment { + pod := NewMultiContainerLogGenerator(namespace, name, 0, delay, message, 1, map[string]string{}) + deployment := runtime.NewDeployment(namespace, name) + labels := map[string]string{ + constants.LabelK8sComponent: "log-generator", + } + deployment.Spec = appsv1.DeploymentSpec{ + Replicas: utils.GetPtr(replicas), + Strategy: appsv1.DeploymentStrategy{ + Type: appsv1.RecreateDeploymentStrategyType, + }, + Selector: &metav1.LabelSelector{ + MatchLabels: labels, + }, + Template: corev1.PodTemplateSpec{ + Spec: pod.Spec, + ObjectMeta: metav1.ObjectMeta{ + Labels: labels, + }, + }, + } + deployment.Spec.Template.Spec.RestartPolicy = corev1.RestartPolicyAlways + return deployment +} + // NewLogGenerator creates a pod that will print `count` lines to stdout, waiting for // `delay` between each line. Lines are of the form " [n] `message`" // where n is the number of lines output so far. Once done printing the pod will diff --git a/test/runtime/observability/cluster_log_forwarder.go b/test/runtime/observability/cluster_log_forwarder.go index aa1e5f6f79..987da06d59 100644 --- a/test/runtime/observability/cluster_log_forwarder.go +++ b/test/runtime/observability/cluster_log_forwarder.go @@ -33,12 +33,16 @@ type PipelineBuilder struct { inputRefs []string } +type ClusterLogForwarderVisitor func(clf *obs.ClusterLogForwarder) type InputSpecVisitor func(spec *obs.InputSpec) type OutputSpecVisitor func(spec *obs.OutputSpec) type FilterSpecVisitor func(spec *obs.FilterSpec) type PipelineSpecVisitor func(spec *obs.PipelineSpec) -func NewClusterLogForwarderBuilder(clf *obs.ClusterLogForwarder) *ClusterLogForwarderBuilder { +func NewClusterLogForwarderBuilder(clf *obs.ClusterLogForwarder, visitors ...ClusterLogForwarderVisitor) *ClusterLogForwarderBuilder { + for _, v := range visitors { + v(clf) + } return &ClusterLogForwarderBuilder{ Forwarder: clf, inputSpecs: map[string]*obs.InputSpec{}, @@ -46,6 +50,11 @@ func NewClusterLogForwarderBuilder(clf *obs.ClusterLogForwarder) *ClusterLogForw } } +// End returns the ClusterLogForwarder when done with the builder +func (b *ClusterLogForwarderBuilder) End() *obs.ClusterLogForwarder { + return b.Forwarder +} + func (b *ClusterLogForwarderBuilder) FromInput(inputType obs.InputType, visitors ...InputSpecVisitor) *PipelineBuilder { visitors = append([]InputSpecVisitor{func(spec *obs.InputSpec) { spec.Type = inputType