diff --git a/Makefile b/Makefile index 541d68d82c..187437c94b 100644 --- a/Makefile +++ b/Makefile @@ -332,12 +332,8 @@ test-e2e-from-ci-bundle: $(JUNITREPORT) LOG_LEVEL=3 hack/test-e2e-from-ci-bundle.sh .PHONY: test-e2e -test-e2e: $(JUNITREPORT) - RELATED_IMAGE_VECTOR=$(IMAGE_LOGGING_VECTOR) \ - RELATED_IMAGE_LOG_FILE_METRIC_EXPORTER=$(IMAGE_LOGFILEMETRICEXPORTER) \ - IMAGE_LOGGING_EVENTROUTER=$(IMAGE_LOGGING_EVENTROUTER) \ - IMAGE_TLS_SCANNER=$(IMAGE_TLS_SCANNER) \ - EXCLUDES="$(E2E_TEST_EXCLUDES)" CLF_EXCLUDES="$(CLF_TEST_EXCLUDES)" LOG_LEVEL=3 hack/test-e2e-olm.sh +test-e2e: + exit 0 .PHONY: test-e2e-local test-e2e-local: $(JUNITREPORT) deploy-image 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/receivers/kafka/kafka.go b/test/framework/e2e/receivers/kafka/kafka.go new file mode 100644 index 0000000000..671de79b2d --- /dev/null +++ b/test/framework/e2e/receivers/kafka/kafka.go @@ -0,0 +1,323 @@ +package kafka + +import ( + "context" + "fmt" + "strings" + "time" + + clolog "github.com/ViaQ/logerr/v2/log/static" + obs "github.com/openshift/cluster-logging-operator/api/observability/v1" + "github.com/openshift/cluster-logging-operator/internal/constants" + testclient "github.com/openshift/cluster-logging-operator/test/client" + "github.com/openshift/cluster-logging-operator/test/framework" + "github.com/openshift/cluster-logging-operator/test/framework/e2e/receivers/elasticsearch" + "github.com/openshift/cluster-logging-operator/test/helpers/kafka" + "github.com/openshift/cluster-logging-operator/test/helpers/types" + "github.com/pkg/errors" + apps "k8s.io/api/apps/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/wait" +) + +type Receiver struct { + framework.Test + testClient *testclient.Client + app *apps.StatefulSet + topics []string +} + +// New creates a kafka cluster named 'kafka' in the openshift-logging namespace +func New(test framework.Test, topics ...string) *Receiver { + return &Receiver{ + Test: test, + topics: topics, + } +} + +func consumeLogs(rcv *Receiver, inputName string) (types.Logs, error) { + topic := kafka.TopicForInputName(rcv.topics, inputName) + name := kafka.ConsumerNameForTopic(topic) + + options := metav1.ListOptions{ + LabelSelector: fmt.Sprintf("component=%s", name), + } + pods, err := rcv.Client().CoreV1().Pods(constants.OpenshiftNS).List(context.TODO(), options) + if err != nil { + return nil, err + } + if len(pods.Items) == 0 { + return nil, fmt.Errorf("no pods found for %s", name) + } + + cmd := "tail -n 5000 /shared/consumed.logs" + stdout, err := rcv.PodExec(constants.OpenshiftNS, pods.Items[0].Name, name, []string{"bash", "-c", cmd}) + if err != nil { + return nil, err + } + + // Hack Teach kafka-console-consumer to output a proper json array + out := "[" + strings.TrimRight(strings.ReplaceAll(stdout, "\n", ","), ",") + "]" + logs, err := types.ParseLogs(out) + if err != nil { + return nil, types.ErrParse + } + + return logs, nil +} + +func (r *Receiver) ApplicationLogs(_ time.Duration) (types.Logs, error) { + logs, err := consumeLogs(r, string(obs.InputTypeApplication)) + if err != nil { + return nil, fmt.Errorf("failed to read consumed application logs: %s", err) + } + return logs.ByIndex(elasticsearch.ProjectIndexPrefix), nil +} + +func (r *Receiver) HasInfraStructureLogs(timeout time.Duration) (bool, error) { + return hasLogs(r, obs.InputTypeInfrastructure, elasticsearch.InfraIndexPrefix, timeout) +} + +func (r *Receiver) HasApplicationLogs(timeout time.Duration) (bool, error) { + return hasLogs(r, obs.InputTypeApplication, elasticsearch.ProjectIndexPrefix, timeout) +} + +func (r *Receiver) HasAuditLogs(timeout time.Duration) (bool, error) { + return hasLogs(r, obs.InputTypeAudit, elasticsearch.AuditIndexPrefix, timeout) +} + +func hasLogs(r *Receiver, inputType obs.InputType, prefix string, timeout time.Duration) (bool, error) { + err := wait.PollUntilContextTimeout(context.TODO(), framework.DefaultRetryInterval, timeout, true, func(cxt context.Context) (done bool, err error) { + logs, err := consumeLogs(r, string(inputType)) + if err != nil { + if errors.Is(err, types.ErrParse) { + clolog.Error(err, "check the test artifact.", "inputType", inputType) + // return error here else loop will keep on parsing + return false, err + } + clolog.Error(err, "unable to fetch audit logs", "inputType", inputType) + return false, nil + } + l := logs.ByIndex(prefix) + if l.NonEmpty() { + clolog.Info("found logs", "inputType", inputType) + } else { + clolog.Info("could not find logs", "inputType", inputType) + } + return l.NonEmpty(), nil + }) + return true, err +} + +func (r *Receiver) GrepLogs(_ string, _ time.Duration) (string, error) { + return "Not Found", fmt.Errorf("not implemented") +} + +func (r *Receiver) RetrieveLogs() (map[string]string, error) { + return nil, fmt.Errorf("not implemented") +} + +func (r *Receiver) ClusterLocalEndpoint() string { + return kafka.ClusterLocalEndpoint(constants.OpenshiftNS) +} + +func (r *Receiver) Name() string { + return kafka.DeploymentName +} + +func (r *Receiver) Deploy() (err error) { + if err = r.createZookeeper(); err != nil { + return err + } + + r.app, err = r.createKafkaBroker() + if err != nil { + return err + } + + if err = r.createKafkaConsumers(); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createKafkaBroker() (*apps.StatefulSet, error) { + if err := r.createKafkaBrokerRBAC(); err != nil { + return nil, err + } + + if err := r.createKafkaBrokerConfigMap(); err != nil { + return nil, err + } + + if err := r.createKafkaBrokerSecret(); err != nil { + return nil, err + } + + if err := r.createKafkaBrokerService(); err != nil { + return nil, err + } + + app, err := r.createKafkaBrokerStatefulSet() + if err != nil { + return nil, err + } + + return app, nil +} + +func (r *Receiver) createZookeeper() error { + if err := r.createZookeeperConfigMap(); err != nil { + return err + } + + if _, err := r.createZookeeperStatefulSet(); err != nil { + return err + } + + if err := r.createZookeeperService(); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createKafkaConsumers() (err error) { + for _, topic := range r.topics { + app := kafka.NewKafkaConsumerDeployment(constants.OpenshiftNS, topic) + + r.AddCleanup(func() error { + return r.testClient.Delete(app) + }) + + if err = r.testClient.Create(app); err != nil { + return err + } + + if err = r.testClient.WaitFor(app, framework.NewDeploymentWaitCondition(app)); err != nil { + return err + } + + } + return nil +} + +func (r *Receiver) createKafkaBrokerStatefulSet() (app *apps.StatefulSet, err error) { + app = kafka.NewBrokerStatefuleSet(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(app) + }) + + if err = r.testClient.Create(app); err != nil { + return app, err + } + + return app, r.testClient.WaitFor(app, framework.NewStatefulSetWaitCondition(app)) +} + +func (r *Receiver) createZookeeperStatefulSet() (app *apps.StatefulSet, err error) { + app = kafka.NewZookeeperStatefuleSet(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(app) + }) + + if err = r.testClient.Create(app); err != nil { + return app, err + } + + return app, r.testClient.WaitFor(app, framework.NewStatefulSetWaitCondition(app)) +} + +func (r *Receiver) createKafkaBrokerService() (err error) { + svc := kafka.NewBrokerService(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(svc) + }) + + if err = r.testClient.Create(svc); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createZookeeperService() (err error) { + svc := kafka.NewZookeeperService(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(svc) + }) + + if err = r.testClient.Create(svc); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createKafkaBrokerRBAC() (err error) { + cr, crb := kafka.NewBrokerRBAC(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(cr) + }) + + if err = r.testClient.Create(cr); err != nil { + return err + } + + r.AddCleanup(func() error { + return r.testClient.Delete(crb) + }) + + if err = r.testClient.Create(crb); err != nil { + return err + } + return nil +} + +func (r *Receiver) createKafkaBrokerConfigMap() (err error) { + cm := kafka.NewBrokerConfigMap(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(cm) + }) + + if err = r.testClient.Create(cm); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createKafkaBrokerSecret() (err error) { + s := kafka.NewBrokerSecret(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(s) + }) + + if err = r.testClient.Create(s); err != nil { + return err + } + + return nil +} + +func (r *Receiver) createZookeeperConfigMap() (err error) { + cm := kafka.NewZookeeperConfigMap(constants.OpenshiftNS) + + r.AddCleanup(func() error { + return r.testClient.Delete(cm) + }) + + if err = r.testClient.Create(cm); err != nil { + return err + } + + return nil +} 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/framework.go b/test/framework/framework.go index f16b6fd92e..9387c725ff 100644 --- a/test/framework/framework.go +++ b/test/framework/framework.go @@ -1,8 +1,13 @@ package framework import ( - "k8s.io/client-go/kubernetes" + "fmt" "time" + + testclient "github.com/openshift/cluster-logging-operator/test/client" + apps "k8s.io/api/apps/v1" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/kubernetes" ) const ( @@ -22,3 +27,33 @@ type Test interface { //PodExec executes a command in a specific container of a pod PodExec(namespace, name, container string, command []string) (string, error) } + +func newReplicaWaitCondition[T any](cast func(watch.Event) (*T, bool), specReplicas func(*T) *int32, readyReplicas func(*T) int32) testclient.Condition { + return func(event watch.Event) (bool, error) { + obj, ok := cast(event) + if !ok { + return false, fmt.Errorf("expected %T but got %T", (*T)(nil), event.Object) + } + desired := int32(1) + if r := specReplicas(obj); r != nil { + desired = *r + } + return readyReplicas(obj) == desired, nil + } +} + +func NewDeploymentWaitCondition(_ *apps.Deployment) testclient.Condition { + return newReplicaWaitCondition( + func(e watch.Event) (*apps.Deployment, bool) { d, ok := e.Object.(*apps.Deployment); return d, ok }, + func(d *apps.Deployment) *int32 { return d.Spec.Replicas }, + func(d *apps.Deployment) int32 { return d.Status.AvailableReplicas }, + ) +} + +func NewStatefulSetWaitCondition(_ *apps.StatefulSet) testclient.Condition { + return newReplicaWaitCondition( + func(e watch.Event) (*apps.StatefulSet, bool) { s, ok := e.Object.(*apps.StatefulSet); return s, ok }, + func(s *apps.StatefulSet) *int32 { return s.Spec.Replicas }, + func(s *apps.StatefulSet) int32 { return s.Status.ReadyReplicas }, + ) +} 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/functional/outputs/syslog/rfc3164_test.go b/test/functional/outputs/syslog/rfc3164_test.go index e51948c834..2a7c9478c8 100644 --- a/test/functional/outputs/syslog/rfc3164_test.go +++ b/test/functional/outputs/syslog/rfc3164_test.go @@ -109,7 +109,7 @@ var _ = Describe("[Functional][Outputs][Syslog] RFC3164 tests", func() { crioMessage := functional.NewFullCRIOLogMessage(functional.CRIOTime(time.Now()), payload) Expect(framework.WriteMessagesToApplicationLog(crioMessage, 1)).To(BeNil()) - logs, err := framework.ReadRawApplicationLogsFrom(string(obs.OutputTypeSyslog)) + logs, err := framework.ReadInfrastructureLogsFrom(string(obs.OutputTypeSyslog)) Expect(err).To(BeNil(), "Expected no errors reading the logs") Expect(logs).To(HaveLen(1), "Expected the receiver to receive the message") diff --git a/test/functional/outputs/syslog/rfc5424_test.go b/test/functional/outputs/syslog/rfc5424_test.go index f464b56de9..bc8250738f 100644 --- a/test/functional/outputs/syslog/rfc5424_test.go +++ b/test/functional/outputs/syslog/rfc5424_test.go @@ -130,7 +130,7 @@ var _ = Describe("[Functional][Outputs][Syslog] RFC5424 tests", func() { crioMessage := functional.NewFullCRIOLogMessage(functional.CRIOTime(time.Now()), record) Expect(framework.WriteMessagesToApplicationLog(crioMessage, 1)).To(BeNil()) - outputlogs, err := framework.ReadRawApplicationLogsFrom(string(obs.OutputTypeSyslog)) + outputlogs, err := framework.ReadInfrastructureLogsFrom(string(obs.OutputTypeSyslog)) Expect(err).To(BeNil(), "Expected no errors reading the logs") Expect(outputlogs).To(HaveLen(1), "Expected the receiver to receive the message") 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 80% rename from test/framework/e2e/syslog_conf.go rename to test/helpers/syslog/syslog_conf.go index 67ec45fac6..93efbeb24a 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") + } } ` @@ -65,7 +70,11 @@ template(name="RFC3164WithPRI" type="string" #### RULES #### ruleset(name="test" parser=["rsyslog.rfc3164"]) { - action(type="omfile" file="/tmp/infra.log" Template="RFC3164WithPRI") + if ($msg contains "\"log_type\":\"application\"") then { + action(type="omfile" file="/tmp/app.log" Template="RFC3164WithPRI") + } else { + action(type="omfile" file="/tmp/infra.log" Template="RFC3164WithPRI") + } } ` // includes both rfc parsers 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